iT邦幫忙

2026 iThome 鐵人賽

DAY 17
0

D16 講完資料怎麼在 DataFusion 內部流:拉、一批一批、跑在 tokio 上,pipeline breaker 是 async 的等待點。今天緊接著問一個沒答案就不能上線的問題:

同一個 tokio runtime 上跑好幾個 pipeline breaker,它們都想要記憶體,記憶體只有一份,誰該讓?

這問題乍看是資源管理小事,但答錯了會直接爆整個 executor process。DataFusion 給的答案是一個叫 MemoryPool 的 trait,兩種預設實作代表兩種公平觀。今天把它拆開看,順便回扣 Comet 一個真痛點:Spark 那邊有 UnifiedMemoryManager 在管 JVM 記憶體、Rust 這邊又有一個 pool,兩個池怎麼對齊。

先看一條會爆的 SQL

假設要跑這一句:

SELECT l_orderkey, SUM(l_extendedprice)
FROM lineitem JOIN orders ON l_orderkey = o_orderkey
WHERE o_orderdate > '1998-01-01'
GROUP BY l_orderkey
ORDER BY 2 DESC

跑一次 EXPLAIN,會看到大概長這樣的樹:

SortExec
  └── AggregateExec (Final)
        └── HashJoinExec
              ├── build side: DataSourceExec (orders)
              └── probe side: DataSourceExec (lineitem)

D16 講過三種是 pipeline breaker——下游要等上游把資料物化完才能開始拉。這裡剛好三個都有:

  • HashJoinExec 的 build side:整張 orders 的 hash table 要建好,lineitem 才能開始 probe
  • AggregateExec (Final):hash aggregation 要把所有 group 累完才能吐第一筆
  • SortExec:要看完全部資料才知道最大的是哪一筆

三個大戶同時要記憶體。整台機器 heap 只有那麼多,如果都不管,OS 一個 OOM 直接把 process 幹掉——連錯誤訊息都沒有,只在 syslog 留一行 killed

所以問題不是「記憶體不夠怎麼辦」,而是不夠的時候誰該讓、讓多少、由誰來判斷

為什麼需要一個顯式的 pool

寫過 C / Rust 的人第一直覺會是「不是有 allocator 嗎,malloc / new 出錯就處理啊」。可是 allocator 只知道「有沒有那麼多 bytes 給你」,不知道「這個 700 MB 是 SortExec 想繼續讀,還是 HashJoin 想蓋 hash table」,也不知道「誰有 spill 路徑、誰沒有」。

要協調,只能有一個更高一層的東西集中看——這就是 MemoryPool

  • 每個要吃記憶體的運算子先跟 pool 註冊一個 MemoryConsumer
  • 用的時候呼叫 reserve(size) 拿一個 MemoryReservation
  • 想再吃更多 grow(additional),吃不完 shrink(released)
  • pool 決定要不要給,超過 limit 就 return 一個 ResourcesExhausted error
  • 拿到 error 的 consumer 自己決定怎麼辦——sort / hash 會去 spill,一般 buffer 就直接 fail

這裡有個很重要的分工:pool 不會自己去 spill 別人。它只做「批准或拒絕」,spill 是 consumer 自己的事。理由是 pool 不知道別人的 state 長怎樣、能不能 spill、spill 到哪裡——這些只有 consumer 自己清楚。

用一個小類比:pool 像公司財務,運算子像各部門。財務只管「本月預算超了沒」,超了就打回申請單。部門收到打回會自己決定要延後案子、還是找別的方案,財務不會幫你決定要砍哪個案子。

這個「批准 / 拒絕」跟「怎麼處理拒絕」分開的設計,是後面 Greedy / Fair 兩種策略能長出來的前提。

兩種公平:Greedy vs Fair

DataFusion 內建兩個主要實作:GreedyMemoryPoolFairSpillPool。名字取得很直白,但取捨要拆開講。

GreedyMemoryPool:先來先贏

規則超簡單:加總所有 reservation,超過 limit 就拒絕新的

先跑的運算子想拿多少拿多少,後跑的一發現額度用完就得馬上 spill。Sort 已經吃了 700 MB、HashJoin 想 grow 200 MB 卻只剩 100 MB 額度——直接拒絕,HashJoin 去 spill,Sort 完全不受影響。

好處是簡單、可預測,適合單 query 或只有一個大戶運算子的場景。壞處是不公平:Sort 明明有 spill 路徑,卻因為它先來就佔盡優勢,HashJoin 明明剛開始蓋 hash table,spill 就會導致後面的 probe 得從 disk 讀。

FairSpillPool:已經吃最多的先讓

規則稍微複雜:把「可 spill 的部分」平均分

  • 每個 consumer 註冊時要說自己是不是 spillable
  • Sort、hash join build 這種自己有 spill 實作的走 spillable 名單
  • 非 spillable 的(例如 RecordBatch 本身)照它 reserve 的量硬扣
  • pool 剩下的額度平分給所有 spillable consumer
  • 有人超過自己那份就被拒絕,變相逼它 spill

關鍵點:Fair 不是「先 spill 誰誰吃虧」,是「已經吃最多的先讓」

回到剛才的例子。假設 pool 上限 1 GB、Sort 已用 700 MB、HashJoin 用 200 MB。新的 grow 進來時:

  • Greedy 看「總和超不超」——超了就拒下一個進來的
  • Fair 看「誰佔比高」——Sort 佔 70% 遠高於 HashJoin 的 20%,超額的份 Sort 該退,HashJoin 反而拿得到

這個切法讓「有 spill 路徑」變成一種義務,而不是運氣——你有 spill 就要有讓的準備,不是你先來就可以霸佔。

兩者怎麼選

沒有絕對答案。單 query 的批次任務用 Greedy 反而穩定,因為 spill 越少越好。多個 spillable 運算子並存(例如同時有 sort + hash join build)Fair 比較合理,不會出現先來的把後來的擠到磁碟。

DataFusion 預設是 GreedyMemoryPool,Comet 那邊會覆蓋這個選擇,等下講。

只追蹤大戶:TrackConsumersPool 的近似

DataFusion 還有第三個工具,叫 TrackConsumersPool——它不是獨立的策略,是包在 Greedy 或 Fair 外面的裝飾器。

功能是:pool 拒絕 reservation 時,錯誤訊息會附上「目前用最多的前 N 個 consumer 各自吃了多少」。這個很實用,OOM 錯誤如果只寫「記憶體滿了」,你根本不知道找誰——是 sort 太貪、還是 hash join 建了太大的表、還是有個 UDF 洩漏。

這個 pool 的近似點在:只維護 top-N 名單,小額配置不追。理由是實務上前幾名通常吃掉八成,追全部要多維護一份 map、每次 reserve / release 都要更新,成本不划算。

問題就出在這個「通常」——不是「總是」。

思考題:什麼場景會讓「只追大戶」這個近似出事?

一個提示:想像一條 query 用了 200 個 filter,每個 filter 各自 buffer 一小塊 batch,單個都排不進 top-10,加起來卻吃掉 300 MB。OOM 時錯誤訊息只會指到那個吃 500 MB 的 sort,但真正該優化的可能是那 200 個小 filter 為什麼要各留一份。

抓錯兇手不會直接害你當機,但會害你花一整個下午在錯的地方調 config。這是近似的隱性成本,發生率不高但一發生就很久。

執行時大概長這樣

把上面拼起來,一條 SQL 的 memory 管理路徑:

      MemoryPool (limit = 1 GB, Fair 策略)
              │
   ┌──────────┼──────────────────────┐
   │          │                      │
 register  register              register
   │          │                      │
   ▼          ▼                      ▼
 SortExec  HashJoinExec           AggregateExec
 (spillable)  build side          (spillable)
              (spillable)
              
   [執行到一半]                  [執行到一半]
   grow(200 MB)                  grow(300 MB)
       │                             │
       ▼                             ▼
   ┌──────────────────────────────────────┐
   │  pool: 我看看                          │
   │  Sort 已 700 MB (70%)                 │
   │  HashJoin 已 200 MB (20%)             │
   │  AggregateExec 已 100 MB (10%)        │
   │  Fair 份額 = 1000 / 3 ≈ 333 MB       │
   │  Sort 超額最多 → 新 grow 拒絕           │
   │  (或選 Sort 那邊拒絕 → 逼 Sort spill)  │
   └──────────────────────────────────────┘
       │
       ▼
   ResourcesExhausted → SortExec 呼叫自己的 spill 路徑
   → 寫一個 run 到 disk → 釋放 reservation → pool 額度回來

幾個實務細節:

  • Reservation 是動態的,一個 batch 內都可能 grow / shrink 好幾次
  • spill 完不代表沒事——後面 sort merge 時還要重新讀 disk,時間成本轉嫁到 I/O
  • 沒有 spill 路徑的 consumer 被拒就是 fail,整條 query 掛
  • pool 是單一 partition 的——多個 partition 共用同一個 pool,還是各自一個,看實作

回到 Comet:兩個池怎麼協調

D15 接點四留了一個問題——JVM 那邊的記憶體歸 Spark 管,Rust 這邊的記憶體歸 DataFusion 管,兩邊怎麼對齊。今天可以正式交代。

Spark 這邊有一個 UnifiedMemoryManager:execution(sort、hash、agg)跟 storage(cache、broadcast)共享一塊 JVM heap,兩塊會動態借用。這是 2015 年 Project Tungsten 定下來的模型,沿用至今。

Rust 這邊有 DataFusion 的 MemoryPool。Comet 每一個 native operator 都跑在這個 pool 底下。

兩者不對齊會出這幾種事:

  • Rust 側覺得還有額度、JVM 側早滿了 → Spark task 觸發 OOM,native 側正在 hash 的資料全丟
  • JVM 側撥了額度給 Comet、Rust 側沒同步 → 兩邊額度加起來超過機器 heap,OS 一個 kill
  • 兩邊各自 spill → Rust 側 spill 到自己的 tmp、JVM 側再 spill 一次同樣的資料,磁碟白寫兩次

Comet 的處理方向是讓 Rust 側的 MemoryPool 變成 Spark manager 的下級

  • native 想 reserve 記憶體時,Comet 這邊的 pool 透過 JNI 呼叫 Spark memory manager 拿額度
  • Off-heap 那塊(Arrow buffer 分配的地方)算在 Spark 的 off-heap pool,避免 Rust 這邊自己搞一份 Spark 看不到的浮空額度
  • Spill 有明確歸屬:Comet 自己的 shuffle spill 走 native 側檔案系統,但額度扣在 Spark 那邊

這條路的代價是每個 reservation 都 pay 一次 JNI 呼叫。D16 講 block_on 是每個 batch 一次來回,那個頻率已經不低;MemoryPool 的 reservation 可能一個 batch 內 grow / shrink 好幾次——JNI 頻率更高。

Comet 從 0.11 開始花力氣減少這個往返(例如 batch reservation、預先撥額度),但結構性成本不會消失,只能攤平。這件事會在 D26 深挖 JNI 邊界時再回來看細節。

結論

MemoryPool 是 DataFusion 把「記憶體協調」抽出來的一層:不做 spill 決策,只做批准 / 拒絕,spill 的怎麼辦交給 consumer。Greedy 跟 Fair 兩種策略代表兩種公平觀,TrackConsumersPool 是包在外面的追蹤器、只留大戶名單當近似。Comet 這邊要把 native pool 對齊到 Spark 的 UnifiedMemoryManager,代價是每個 reservation 都要跨 JNI 問一次。

明天 D18 接執行細節的第三題:Join 演算法怎麼選?hash、merge、symmetric hash、nested loop 各自的取捨在哪,planner 憑什麼在這幾個裡挑一個?其中一個很硬的因素就是今天的主題——這個演算法要多少記憶體、有沒有 spill 路徑。

那就明天見~

參考資料


上一篇
Day 16 DataFusion 執行模型:Volcano Pull、RecordBatch 與 Tokio Async
下一篇
Day 18 Join 演算法:planner 憑什麼選
系列文
1+1+1>3 ~ Spark 與 DataFusion、Comet 效能煉金術 ~21
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言